Skip to content

无锁环形队列设计与实现(SPSC 与 MPMC) ​

生产者—消费者是并发程序中最常见的协作模式之一。用互斥锁保护队列,在竞争激烈时会带来线程阻塞与上下文切换的开销;无锁环形队列改用原子操作和内存序保证线程安全,读端与写端互不阻塞。本文从最简单的单生产者单消费者(SPSC)环形队列讲起,说明空满区分与内存序的配对方式;再分析多生产者多消费者(MPMC)的实现难点,给出 Dmitry Vyukov 提出的按槽序号(sequence)方案的完整实现。适合需要在 C++ 中手写或读懂无锁队列的开发者阅读。

SPSC:数据结构 ​

单生产者单消费者是最简单的并发模型:生产者是写索引的唯一修改者,消费者是读索引的唯一修改者。这一「索引各归一线程」的性质,是 SPSC 队列可以用极少的同步机制实现的前提。

  • 固定容量数组(buffer):存储元素的连续内存空间,大小为 capacity + 1。额外的一个空槽用于区分「满」和「空」两种状态。
  • 读索引(head)与写索引(tail):head 指向下一个要读取的位置,tail 指向下一个要写入的位置,均为 std::atomic<size_t>。
  • 循环递增:索引按模 capacity + 1 递增,实现缓冲区的循环利用。

用「空一格」方案区分空满:队列满的条件是 (tail + 1) % (capacity + 1) == head,队列空的条件是 head == tail。如果不用这个额外空槽,head == tail 将同时表示空和满,需要额外的标志位或计数器才能区分。

SPSC:操作流程与内存序 ​

写入(push):

  1. 读取当前 tail;
  2. 计算 next_tail = (tail + 1) % (capacity + 1);
  3. 若 next_tail == head,队列已满,返回失败;
  4. 将数据写入 buffer[tail];
  5. 以 memory_order_release 更新 tail = next_tail。

读取(pop):

  1. 读取当前 head;
  2. 若 head == tail,队列为空,返回失败;
  3. 读取 buffer[head];
  4. 以 memory_order_release 更新 head = (head + 1) % (capacity + 1)。

两个方向的 release/acquire 配对,是正确性的核心:

  • tail 方向:生产者写完数据后 tail.store(release);消费者 tail.load(acquire) 读到新 tail 时,保证能看到对应槽位上已完整写入的数据。
  • head 方向:消费者读完数据后 head.store(release);生产者 head.load(acquire) 读到新 head 后才复用该槽位,保证覆盖前消费者已经读完旧数据。

反过来,读写自己的索引(生产者读 tail、消费者读 head)用 memory_order_relaxed 即可:该索引只有一个写者,不存在与其他线程同步的问题。如果索引完全不用原子类型,读写索引构成数据竞争,在 C++ 标准下是未定义行为,可能出现撕裂读或索引更新丢失。

SPSC:完整实现 ​

cpp
#include <atomic>
#include <vector>
#include <optional>
#include <cassert>

template<typename T>
class LockFreeRingBuffer {
public:
    explicit LockFreeRingBuffer(size_t capacity)
        : capacity_(capacity + 1),  // 一个额外空间用来区分满和空
          buffer_(capacity_),
          head_(0),
          tail_(0) {
        assert(capacity > 0);
    }

    // 禁止复制
    LockFreeRingBuffer(const LockFreeRingBuffer&) = delete;
    LockFreeRingBuffer& operator=(const LockFreeRingBuffer&) = delete;

    // 生产者调用:插入元素,成功返回true,满返回false
    bool push(const T& item) {
        size_t tail = tail_.load(std::memory_order_relaxed);
        size_t next_tail = increment(tail);

        if (next_tail == head_.load(std::memory_order_acquire)) {
            // 缓冲区满
            return false;
        }

        buffer_[tail] = item;
        tail_.store(next_tail, std::memory_order_release);
        return true;
    }

    // 生产者移动语义版本
    bool push(T&& item) {
        size_t tail = tail_.load(std::memory_order_relaxed);
        size_t next_tail = increment(tail);

        if (next_tail == head_.load(std::memory_order_acquire)) {
            return false;
        }

        buffer_[tail] = std::move(item);
        tail_.store(next_tail, std::memory_order_release);
        return true;
    }

    // 消费者调用:取出元素,存在时返回 optional<T>,无数据时返回 empty optional
    std::optional<T> pop() {
        size_t head = head_.load(std::memory_order_relaxed);
        if (head == tail_.load(std::memory_order_acquire)) {
            // 缓冲区空
            return std::nullopt;
        }

        T item = std::move(buffer_[head]);
        head_.store(increment(head), std::memory_order_release);
        return item;
    }

    // 是否空
    bool empty() const {
        return head_.load(std::memory_order_acquire) == tail_.load(std::memory_order_acquire);
    }

    // 容量(实际可用容量为 capacity_ - 1)
    size_t capacity() const {
        return capacity_ - 1;
    }

private:
    size_t increment(size_t idx) const noexcept {
        return (idx + 1) % capacity_;
    }

private:
    const size_t capacity_;
    std::vector<T> buffer_;
    std::atomic<size_t> head_;
    std::atomic<size_t> tail_;
};

为什么 MPMC 复杂得多 ​

一旦存在多个生产者或多个消费者,「索引各归一线程」不再成立:

  • 同一索引被并发修改:多个生产者同时推进 tail,必须用 CAS(Compare-And-Swap)循环保证推进操作的原子与串行化,失败方重试。
  • ABA 问题:一个值从 A 改为 B 又改回 A,依赖「值未变化」判断的 CAS 会误判,导致同一槽位被重复占用或覆盖。
  • 槽位状态管理:多个生产者必须保证每次写入独占一个槽位;多个消费者必须保证每个元素只被消费一次。空、满、写入中、可读等状态需要额外的标记。
  • 虚假共享(false sharing):多个线程频繁读写落在同一 CPU 缓存行上的计数器,会导致缓存行反复失效,显著影响性能。

链表型的经典解法是 Michael-Scott 队列(1996),但它引入了已出队节点的内存回收问题,需要延迟回收机制配合。环形数组方案则把问题转化为「每个槽位的状态管理」,回避了动态内存回收。

MPMC:Vyukov 的按槽序号方案 ​

Dmitry Vyukov 在 1024cores.net 发表的 bounded MPMC queue 给出了一种被广泛借鉴的设计,核心是给每个槽位一个独立序号:

  1. 初始化:第 i 个槽的 sequence 为 i;两个全局位置计数 enqueue_pos_(生产者推进)与 dequeue_pos_(消费者推进)从 0 开始。
  2. 生产者抢占:要写入位置 pos 时,检查槽序号 seq 与 pos 的差:
    • diff == 0:该槽在当前这一圈空闲,用 CAS 把 enqueue_pos_ 从 pos 抢到 pos + 1,成功者独占该槽;
    • diff < 0:消费者还没释放该槽(序号落后于位置),队列已满;
    • diff > 0:其他生产者已经抢先推进了位置,重新加载 enqueue_pos_ 再试。
  3. 写入与发布分离:抢占成功后写入数据,再把该槽 sequence 发布为 pos + 1——消费者只有在看到 seq == pos + 1 时才认为数据有效。
  4. 消费者抢占:要读取位置 pos 时,检查 seq == pos + 1 表示数据有效,同样用 CAS 抢 dequeue_pos_;读完后把 sequence 置为 pos + capacity,即该槽在下一圈的位置,对未来的生产者重新可写。

这个设计的两个关键点:

  • CAS 只抢位置计数,不携带数据。CAS 失败的线程只需重载位置重试;数据写入在独占槽位之后单独进行,用 release 发布序号。
  • 序号天然规避 ABA:sequence 等于「位置 + 已完成圈数 × capacity」,随圈数单调递增,同一槽位每一圈的序号都不同,旧状态不会被误认为当前状态。

需要说明的是,该算法严格意义上并非 lock-free:如果一个线程在抢占槽位之后、发布数据之前被长期挂起,等待该位置数据的消费者会停住。它是不使用互斥锁的非阻塞实现,工程语境中通常仍归入无锁队列;Vyukov 本人在原文中也注明了这一点。

MPMC:完整实现 ​

cpp
#include <atomic>
#include <vector>
#include <cassert>
#include <optional>

template <typename T>
class MPMCRingBuffer {
public:
    explicit MPMCRingBuffer(size_t capacity)
        : buffer_(capacity), capacity_(capacity), index_mask_(capacity - 1)
    {
        // capacity必须是2的幂,这样方便用掩码计算index
        assert((capacity >= 2) && ((capacity & (capacity - 1)) == 0));

        for (size_t i = 0; i < capacity; ++i) {
            buffer_[i].sequence.store(i, std::memory_order_relaxed);
        }

        enqueue_pos_.store(0, std::memory_order_relaxed);
        dequeue_pos_.store(0, std::memory_order_relaxed);
    }

    // 禁止复制
    MPMCRingBuffer(const MPMCRingBuffer&) = delete;
    MPMCRingBuffer& operator=(const MPMCRingBuffer&) = delete;

    bool push(const T& data) {
        Cell* cell;
        size_t pos = enqueue_pos_.load(std::memory_order_relaxed);

        for (;;) {
            cell = &buffer_[pos & index_mask_];
            size_t seq = cell->sequence.load(std::memory_order_acquire);
            intptr_t diff = (intptr_t)seq - (intptr_t)pos;

            if (diff == 0) {
                if (enqueue_pos_.compare_exchange_weak(pos, pos + 1, std::memory_order_relaxed)) {
                    break;
                }
            } else if (diff < 0) {
                // 队列满
                return false;
            } else {
                // 重新加载生产者位置,准备重试
                pos = enqueue_pos_.load(std::memory_order_relaxed);
            }
        }

        cell->data = data;
        cell->sequence.store(pos + 1, std::memory_order_release);
        return true;
    }

    std::optional<T> pop() {
        Cell* cell;
        size_t pos = dequeue_pos_.load(std::memory_order_relaxed);

        for (;;) {
            cell = &buffer_[pos & index_mask_];
            size_t seq = cell->sequence.load(std::memory_order_acquire);
            intptr_t diff = (intptr_t)seq - (intptr_t)(pos + 1);

            if (diff == 0) {
                if (dequeue_pos_.compare_exchange_weak(pos, pos + 1, std::memory_order_relaxed)) {
                    break;
                }
            } else if (diff < 0) {
                // 队列空
                return std::nullopt;
            } else {
                pos = dequeue_pos_.load(std::memory_order_relaxed);
            }
        }

        T result = cell->data;
        cell->sequence.store(pos + capacity_, std::memory_order_release);
        return result;
    }

private:
    struct Cell {
        std::atomic<size_t> sequence; // 序号,用于标志和版本控制槽状态
        T data;                       // 数据存储
    };

    std::vector<Cell> buffer_;         // 环形数组数据槽
    size_t const capacity_;            // 容量(2的幂)
    size_t const index_mask_;          // 用于快速取模(capacity_ - 1)

    std::atomic<size_t> enqueue_pos_; // 生产者索引
    std::atomic<size_t> dequeue_pos_; // 消费者索引
};

实现细节上:容量要求为 2 的幂,这样 pos & index_mask_ 等价于取模且更快;槽的 sequence 用 acquire 读、release 写,与数据写入/读取构成发布—获取配对;生产者读取自己的 enqueue_pos_ 用 relaxed,因为 CAS 本身会仲裁位置归属。

要点 ​

  • SPSC 的最小正确实现只需要两个原子索引、一个额外空槽和两组 release/acquire 配对;能用 SPSC 的场景不必引入更复杂的机制。
  • MPMC 的关键在于按槽独立序号:CAS 只用于抢占位置计数,数据写入与序号发布分离,序号随圈数单调递增从而规避 ABA。
  • 「无锁」要区分语境:Vyukov 方案严格意义上不是 lock-free,但互不阻塞、无互斥锁,实践中满足多数高性能场景。
  • 优先成熟库:无锁结构调试困难、内存序容易出错。boost::lockfree::queue 与 boost::lockfree::spsc_queue、folly::MPMCQueue、moodycamel::ConcurrentQueue 等实现经过充分验证,生产环境通常应先于手写实现考虑。
最近更新

基于 VitePress 构建